stream: clean up when pipeline throws synchronously - #65064
Conversation
|
Review requested:
|
9c3b6a8 to
cb6d899
Compare
ronag
left a comment
There was a problem hiding this comment.
I would say that ownership is not taken until pipeline succeeds...
There was a problem hiding this comment.
Pull request overview
This PR fixes a resource-leak in stream.pipeline() when pipelineImpl() throws synchronously during its wiring loop (e.g. ERR_STREAM_UNABLE_TO_PIPE), ensuring already-adopted streams are torn down and abort listeners are removed before the error propagates.
Changes:
- Wrap the stream-wiring loop in
pipelineImpl()with synchronous-throw cleanup that drainsdestroys, disposes the outerAbortSignallistener, and aborts the internal controller. - Add a regression test covering synchronous-throw cleanup for
ERR_STREAM_UNABLE_TO_PIPEandERR_INVALID_RETURN_VALUE, including abort-listener disposal in the promise form.
Reviewed changes
Copilot reviewed 2 out of 2 changed files in this pull request and generated 1 comment.
| File | Description |
|---|---|
| lib/internal/streams/pipeline.js | Adds teardown on synchronous throws during the wiring loop to prevent leaked streams/fds and leaked abort listeners. |
| test/parallel/test-stream-pipeline-sync-throw-cleanup.js | New regression test asserting cleanup happens on synchronous throw paths (fd release, stream destroyed, abort listener removed). |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| } catch (err) { | ||
| // The loop above can throw synchronously (e.g. ERR_STREAM_UNABLE_TO_PIPE) | ||
| // after some streams have already been wired up. Those streams are | ||
| // registered in `destroys` but `finishImpl()` never runs, so tear them | ||
| // down here before propagating, otherwise their resources leak. | ||
| while (destroys.length) { | ||
| destroys.shift()(err); | ||
| } | ||
| disposable?.[SymbolDispose](); | ||
| ac.abort(); | ||
| throw err; | ||
| } |
|
Thanks @ronag, that is a fair point and I think it splits the PR into two separate questions. Destroying the streams. I accept the ownership argument. If The So even setting the destroy question aside, Would you prefer I narrow this to only undoing One doc question while we are here: On the Copilot comment about |
`pipelineImpl()` adds a listener to the caller's `AbortSignal` before it wires the streams together, and only disposes of it in `finishImpl()`. The wiring loop can throw synchronously, for example `ERR_STREAM_UNABLE_TO_PIPE` when the destination is already destroyed, and `finishImpl()` never runs on that path, so the listener stays attached for the lifetime of the signal. A long lived signal reused across many failed calls accumulates one listener per call. Dispose of it before propagating the error. The streams themselves are left untouched, since ownership is not taken until the pipeline has been established, which is the behaviour the existing `ERR_INVALID_RETURN_VALUE` cases in `test-stream-pipeline.js` assert. Signed-off-by: Shani Singh <teamdeveloperworld@gmail.com>
cb6d899 to
c4f4c9f
Compare
|
You were right, and the test suite says so explicitly. CI caught it: That line is I have narrowed this to only the part that is not about ownership: disposing the listener that I also re-ran the three The test in this PR now only covers the signal listener, since the not-destroyed behaviour is already covered by the existing cases. |
pipelineImpl()wires the streams together in a loop that can throw synchronously. The most common case isERR_STREAM_UNABLE_TO_PIPE, raised when the next stream is already closed or destroyed, which happens routinely when a destination goes away first (for examplepipeline(fs.createReadStream(file), res)after the HTTP client disconnected).Each stream the loop adopts registers a destroy function in
destroys.finishImpl()is the only code that drainsdestroys, disposes the listener added to the caller'sAbortSignaland callsac.abort(), and it never runs when the loop throws. Every stream already wired up is therefore left undestroyed and its resources leak. For anfs.ReadStreamsource that is a leaked file descriptor.The loop has six synchronous throw sites: one
ERR_STREAM_UNABLE_TO_PIPE, threeERR_INVALID_RETURN_VALUEand twoERR_INVALID_ARG_TYPE. This wraps the loop so the same teardown runs before the error propagates. The error is still thrown, so the observable failure mode is unchanged.Most of the diff is the re-indentation of the existing loop. Reviewing with
git diff -wshows the actual change is 13 lines.Before
After
The full reproduction is in the linked issue. The added test fails on
main(5 failing assertions) and passes with this change.I also checked the change against a set of ordinary
pipeline()usages (happy path,fsread to writable, asynchronous mid-stream error,ENOENTsource, async generator transform, abort via an outer signal) and the behaviour is identical before and after. The only observable differences are on the two synchronous-throw paths, where the streams are now destroyed.Fixes: #65063